MQ - RocketMQ 5.4 集群模式安装和使用
RocketMQ 集群安装说明
从单机走向集群,是任何分布式系统迈向工业生产线的必经之路。在生产环境中,单机模式就像在钢丝绳上跳舞,任何一次突发的断电、硬件故障、磁盘损坏,都可能会导致业务系统全线瘫痪。RocketMQ 集群模式最核心的设计点可以总结为两句话:
- NameServer 无状态对等:NameServer 之间互不通信,各自分立。Broker 启动时会向所有 NameServer 报到。这种 “无状态” 设计意味着即使挂掉大半,只要剩下一台,整个集群的路由大脑就依然清醒。
- Broker 主从双写与水平分片:
- 水平分片(分担压力):不同的 Master 节点(如 Master A 和 Master B)各自承载 Topic 的一部分队列,多 Master 并存可以物理叠加集群的写吞吐量。
- 主从多副本(防死防丢):每个 Master 绑定一个 Slave。Master 负责读写,Slave 负责实时同步备份。当 Master 挂掉,消费者会自动平滑切换到 Slave 读取数据,确保业务不流断。
在生产落地时,通常有以下三种架构路线:
- 多 Master 模式:全是 Master,没有 Slave。优点是配置简单,缺点是一旦某台机器炸了,未消费的消息在机器恢复前物理不可访问。
- 多 Master 多 Slave 模式(异步复制):性能高,但 Master 瞬间断电,可能存在毫秒级的数据未同步丢失。
- 多 Master 多 Slave 模式(同步双写 / DLedger 纯Raft):金融生产级标准。消息必须强行落盘到 Slave 成功后才返回给客户端,数据零丢失,且 Raft 协议支持 Master 挂掉后 Slave 自动竞选提拔为新 Master。
为了兼顾高性能与高可用,我们今天采用业界应用最广、最稳健的经典模型:2主2从(2M-2S)异步复制、同步刷盘集群架构。我们准备两台物理机(或者虚拟机),通过交叉部署 Master 和 Slave,在节约机器成本的同时,实现硬件级容灾。如果 192.168.1.7 整台机器彻底烧毁,由于 broker-a 的副本 broker-a-s 活着在 .8 上,broker-b 的主节点也活着在 .8 上,整个集群的业务数据依然全面可访问。

服务端集群的安装步骤
单机版的安装可以直接参考 RocketMQ 5.x 的安装和使用。这里我们直接切入最核心的集群配置文件编写。RocketMQ 官方在 conf/ 目录下默认提供了 2m-2s-async/ 模板,我们直接对其进行改造。
机器1的配置
在 .7 机器上,我们需要配置两个文件,一个负责 Master A,一个负责 Slave B。
文件一:配置 Master A。创建并编辑 conf/2m-2s-async/broker-a.properties:
1 | # 集群所属名称 |
文件二:配置 Slave B。创建并编辑 conf/2m-2s-async/broker-b-s.properties:
1 | brokerClusterName=DefaultCluster |
机器2的配置
在 .8 机器上,同样配置两个文件,负责 Master B 和 Slave A。
文件一:配置 Master B。创建并编辑 conf/2m-2s-async/broker-b.properties:
1 | brokerClusterName=DefaultCluster |
文件二:配置 Slave A。创建并编辑 conf/2m-2s-async/broker-a-s.properties:
1 | brokerClusterName=DefaultCluster |
集群启动与验证
配置完成后,严格按照以下时序在两台机器上启动对应的系统组件。
步骤 1:两台机器同时拉起 NameServer。在 192.168.1.7 和 192.168.1.8 上分别执行:
1 | $ nohup sh bin/mqnamesrv & |
步骤 2:在 192.168.1.7 上启动 Broker 矩阵
1 | # 1. 启动 Master A |
步骤 3:在 192.168.1.8 上启动 Broker 矩阵
1 | # 1. 启动 Master B |
步骤 4:用命令检验集群成果
1 | sh bin/mqadmin clusterList -n 192.168.1.7:9876 |
客户端的使用
客户端代码
集群搭建完毕后,我们在之前单机版的 Spring Boot / Java 项目中接入时,代码改动极小。核心的物理变化是:NameServer 的地址可以改为传入全部的节点列表(当然也可以传一个或者一部分)。
1 | import org.apache.rocketmq.client.producer.DefaultMQProducer; |
遇到的问题
当这行代码执行时,Producer 会随机挑一台 NameServer(如 .7)拉取 TOPIC_ORDER 的路由表。由于我们配置了 2 个 Master,NameServer 会返回告诉它:这个 Topic 目前在 broker-a 上有 4 个队列,在 broker-b 上也有 4 个队列。随后,你的客户端在发消息时,就会自动轮询向 192.168.1.5:10911 (Master A) 和 192.168.1.7:10911(Master B) 交叉砸入字节流。整个系统的单点瓶颈得到完美解决。好,我们来运行上述代码,看看消息发送时是不是真的轮询往 MasterA 和 MasterB 发呢?
1 | 消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C139BE22D8CFE0241284B40000, offsetMsgId=C0A8010700002A9F0000000000038BF7, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=0], queueOffset=0, recallHandle=null] |
嗯?奇怪!怎么不往 broker-b 发呢?而且无论我们如何在客户端发消息,就是出不来 broker-b!😓
问题的解决
其实在集群模式下,导致这种现象的原因只有一个:因为你在代码里,没有为 TOPIC_ORDER 提前创建路由元数据,导致 broker-b 根本没有向 NameServer 认领这个 Topic.。罪魁祸首就是我们在 Broker 的配置文件里都设置了 autoCreateTopicEnable=true。
- 当你的客户端发出第一条消息时,它去 NameServer 问:“TOPIC_ORDER 在哪?”
- 此时因为是新 Topic,NameServer 账本里根本没有它!NameServer 就会把系统默认的一个兜底模板 Topic(TBW102)的路由丢给客户端。
- TBW102 默认会指向当前第一个连接最快、最健康的 Broker(也就是刚才的 broker-a)。
- 于是,客户端就在本地内存里为 TOPIC_ORDER 绑定了 broker-a 的 4 个队列(Queue 0, 1, 2, 3)。这就是为什么我们刚才的日志里只出现了 broker-a 的 2 -> 3 -> 0 -> 1 的轮询。
- 既然客户端在内存里已经认定 TOPIC_ORDER 属于 broker-a,那么不管你 sleep 500ms 还是 5000ms,客户端每 30 秒去跟 NameServer 同步路由时,NameServer 返回的依然只有 broker-a。
- 只有当有人真正往 broker-b 发送一次该 Topic 的消息,broker-b 才会动态在磁盘上为它创建文件夹并上报 NameServer。这就陷入了一个 “broker-b 不建立路由 $\rightarrow$ 客户端不往它发 $\rightarrow$ 它更不建立路由” 的死循环。
在真正的生产环境里是绝对禁止开启 autoCreateTopicEnable=true 的。所有的 Topic 必须在部署时由运维人员通过命令手动在所有 Master 上整齐划一地创建好。现在,我们只能在 centos10-01 上执行以下命令,强行打通 broker-b 的任督二脉:
1 | # config/topics.json |
再次重启你的 Java 业务程序,客户端这次再向 NameServer 询问。这一次 NameServer 会吐出一个完美的、包含了 broker-a(4个队列)和 broker-b(4个队列)共计 8 个队列的完整路由表。
1 | 消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C13B7D22D8CFE0241FBA6C0005, offsetMsgId=C0A8010700002A9F000000000004768A, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=0], queueOffset=12, recallHandle=null] |
sendAsync 和 send 的执行
异步发送 sendAsync
异步发送利用了 TCP 通道复用。业务线程把字节流往 Netty 管道里一扔,瞬间返回(耗时接近 0ms)。当 Broker 的响应顺着网线飞回来时,Netty 读线程会根据 opaque 流水号去认亲,并执行对应的 SendCallback。通常用于高并发、大吞吐、允许轻微延迟的非核心主线业务(如日志收集、用户行为埋点、给用户发优惠券通知)。异步发送最突出作用就是能够释放网卡吞吐量。
1 | import org.apache.rocketmq.client.producer.DefaultMQProducer; |
1 | 【异步发送】开始投递... |
在 RocketMQ 的物理世界里,客户端和 broker-a 之间,通常只有一条(或者极少数几条)长连接 TCP 管道。当 100 个业务线程并发调用 sendAsync 时,这些消息就像是 100 颗不同颜色的子弹,被无锁(或者极轻量锁)地倾泻进同一个 TCP 管道中。由于网线是串行的,这 100 颗子弹在网线上是相互交错、前后挨着的。 Broker 收到后,由于 netty 采用了 SendMessageExecutor 业务线程池并发写盘,写盘有快有慢。 这就导致 Broker 吐回来的响应顺序,跟客户端发过去的请求顺序,是完全对不上的!比如客户端先发了请求 1,再发了请求 2。但 Broker 可能会先返回请求 2 的结果,再返回请求 1 的结果。既然网络全乱了,客户端怎么知道哪一个响应属于哪一个业务回调? 答案就是核心组件:opaque(请求唯一标识)与 ResponseFuture。
在 RocketMQ 的通讯内核 rocketmq-remoting 中,有一个核心类叫 ResponseFuture。每次你调用异步发送,底层都会在内存中凭空组装出这样一个精密的状态机:
1 | public class ResponseFuture { |
我们来看底层的执行代码到底是怎么走的:
第一,业务线程挂载账本(不阻塞):
- 你的业务线程调用 sendAsync。
- RocketMQ 底层通过一个内部原子计数器(AtomicInteger),生成一个全局唯一的 opaque(比如:9527)。
- 把你的 SendCallback 业务回调函数包装进 ResponseFuture 对象中。
- 关键一步是把这个对象塞进一个客户端全局的内存哈希表里:responseTable.put(9527, future)。这个表就是“认亲账本”。
- Netty 异步把消息通过 TCP 网线推出去。业务线程完全不等待,立刻挥一挥衣袖返回。线程重获自由,可以去处理下一个请求。
第二,网线与 Broker 端,疯狂并网:
- 请求夹带着 opaque = 9527 的标记划过网线。
- Broker 收到后,高并发写盘。处理完毕后,组装响应报文,并原封不动地把 opaque = 9527 塞回响应头里,顺着网线吐回给客户端。
第三,Netty 读线程(Selector),顺藤摸瓜,惊醒回调:
- 客户端的网卡收到了二进制字节流,Netty 的 Selector 线程(I/O 线程)抓起这串字节,经过编解码器翻译,还原成一个 RemotingCommand 响应对象。
- Netty 读线程一眼看到了这个响应对象里面的 opaque = 9527。
- 它立刻冲向刚才那个 “认亲账本”:ResponseFuture future = responseTable.remove(9527)。
- 拿到这个 future 后,Netty 读线程(或者丢给公用的异步回调线程池)直接调用 future.getInvokeCallback().operationComplete(future)
- 顺着这个入口,最终点燃了你在 Java 业务层里写下的 onSuccess(SendResult sendResult) 代码块!
RocketMQ 压榨网卡的终极魔法,就是用一个基于全局唯一流水号(opaque)的内存认亲表(responseTable),把串行的 TCP 网线强行扭置成了并发的立交桥。网络 I/O 线程只管收发,业务线程只管扔消息,彼此通过 ResponseFuture 隔空对齐。
同步发送 send
同步发送会物理挂起当前业务线程(利用 CountDownLatch),直到 Broker 返回 SendResult。通常用于对数据正确性要求极高、不能容忍任何丢失的场景(如钱包扣款、订单创建)。
1 | import org.apache.rocketmq.client.producer.DefaultMQProducer; |
1 | 【同步发送】准备砸入网络,当前业务线程即将挂起休眠... |
如果你调用的是同步的 producer.send(msg),它的底层网络架构和 sendAsync 一模一样,用的也是同一个 TCP 管道,也生成了 opaque。唯一的物理区别在于:
同步发送时,业务线程把消息交给 Netty 后,没有立刻返回,而是抓起 ResponseFuture 内部的那个 countDownLatch,调用了 future.countDownLatch.await(timeout, TimeUnit.MILLISECONDS)
业务线程原地物理挂起(休眠)。
当 Netty 读线程根据 opaque = 9527 捞出 future 时,它不会执行回调,而是执行
future.countDownLatch.countDown()
这一脚油门下去,把正在休眠的业务线程瞬间唤醒,业务线程从 future.responseCommand 里拿走结果,假装自己完成了一次 “同步阻塞” 调用。
消息接收端代码注意事项
在 RocketMQ 的物理世界里,消息接收端(Consumer)通常是整个微服务分布式架构中最脆弱、最容易沦陷的一环。因为发送端(Producer)往往只负责高频地往外抛字节,而接收端却要实打实地承载执行业务逻辑、写数据库、调第三方接口等沉重的肉身。如果接收端代码写得不规范,极易引发内存溢出(OOM)、消息积压、顺序错乱、以及死锁等生产事故。
必须死守幂等性
消息接收端必须死守“幂等性”,绝不盲目信任中间件。由于网络抖动、重试机制或消费者扩容重平衡,RocketMQ 只能保证 “At least once”(至少投递一次),绝对无法避免重复投递。
❌ 错误示范:拿到消息直接做业务
1 | // 错误:一旦网络抖动,这条消息被投递了两次,用户就会被扣款两次! |
✅ 正确示范:利用数据库唯一键或 Redis 去重表实施物理拦截
1 | consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> { |
区分并发消费与顺序消费
RocketMQ 提供了两种监听器:
- MessageListenerConcurrently(并发消费):多线程同时吃不同的消息,速度极快,但无法保证顺序。
- MessageListenerOrderly(顺序消费):单线程锁定同一个 Queue,严格按照物理写入顺序一条一条吃。如果你的业务(如订单状态流转:创建 $\rightarrow$ 支付 $\rightarrow$ 发货)要求绝对不能乱,就必须使用 Orderly。
❌ 错误示范:用并发监听器去处理顺序业务
1 | // 错误:虽然你用 MessageQueueSelector 投递到了同一个 Queue,但在这里被多线程并发消费, |
✅ 正确示范:使用 MessageListenerOrderly 死守单通道
1 | // 正确:RocketMQ 底层会自动用线程锁、内存锁把该队列死死锁住,强迫在单线程下顺次消费 |
严格禁止在监听器内将消息抛给自定义线程池处理
我们已经知道,消费点位(Offset)的右移,完全取决于消费监听器返回 CONSUME_SUCCESS 的时机。如果你在监听器内部自己搞了一个 new ThreadPoolExecutor 把消息丢进去异步执行,然后监听器立刻返回了 SUCCESS,这会导致 RocketMQ 的自愈链路完全失效!
❌ 错误示范:自作聪明的“异步消費”
1 | // 致命错误:消息刚丢给自定义线程池,监听器就返回 SUCCESS 了,Broker 端认为消费成功,Offset 直接右移。 |
✅ 正确示范:老老实实利用 RocketMQ 自带的线程参数调优。如果你觉得消费不够快,正确的姿势是去调大 Consumer 自带的底层线程池参数,而不是自己写线程池:
1 | DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("GID_INDUSTRIAL_CONSUMER"); |
重试和异常情况的处理
并发消费下:捕获一切异常,优雅返回 RECONSUME_LATER。在 Concurrently 模式下,不要尝试在代码里自己写 while(true) 死循环去重试。应该主动返回 RECONSUME_LATER,让消息进入我们第三章拆解的 %RETRY% 梯度阻尼重试队列,让出 CPU 算力。
1 | consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> { |
顺序消费下:千万不能返回 SUSPEND 太久。在 Orderly 顺序模式下,如果返回了 SUSPEND_CURRENT_QUEUE_A_MOMENT(挂起一会再试),RocketMQ 默认会每隔 1 秒钟疯狂重新投递这条失败的消息,直到成功。
- 如果你的底层数据库彻底瘫痪了,代码一直返回 SUSPEND,这条消息会把整个 Queue 文件夹死死卡住几个小时,后面的订单消息全部积压。
- 在顺序消费时,代码内部自己记录重试次数,如果超过 3 次依然失败,记录日志,强行返回 SUCCESS 释放通道,转为人工补偿,切记!